Conversation
1979696 to
4735398
Compare
Codecov Report❌ Patch coverage is Additional details and impacted files@@ Coverage Diff @@
## main #25304 +/- ##
==========================================
- Coverage 81.91% 81.91% -0.01%
==========================================
Files 1134 1134
Lines 425708 425722 +14
Branches 425708 425722 +14
==========================================
+ Hits 348726 348729 +3
Misses 56305 56305
- Partials 20677 20688 +11 ☔ View full report in Codecov by Harness. 🚀 New features to boost your workflow:
|
| if hash_join.inputs_satisfy_partitioned_requirements()? { | ||
| return Ok(None); | ||
| } | ||
|
|
There was a problem hiding this comment.
Thanks, @stuhood. I tested the Hive-partitioned case and found a wrong-results regression with preserve_file_partitions=1.
There was a problem hiding this comment.
This can drop rows with preserve_file_partitions=1: Hive scans advertise Hash([k], n), but different value sets can place matching keys in different partitions (#23436).
With this branch and target_partitions=3:
- Dim {A,B,C,D} → [[A,D],[B],[C]]
- Fact {B,C,D} → [[B],[C],[D]]
Joining on k returns 0 rows, versus 3 with this hunk reverted. The underlying bug already affects non-broadcast joins; this PR extends it to small tables previously protected by CollectLeft.
Dynamic-filter pushdown already rejects Hash/Hash under this setting, which also explains the lost filters in the updated plans. Could we apply the same guard here and add this regression test?
There was a problem hiding this comment.
I believe that that is a bug in another portion of the codebase (as you pointed out): #25279 is out in order to fix it.
There was a problem hiding this comment.
Thanks for pointing me to #25279, I missed that existing fix. That looks like the right place to address this.
Which issue does this PR close?
Rationale for this change
In
JoinSelection,try_collect_leftconverts a hash join toPartitionMode::CollectLeft(broadcast) when either input is below configured byte/row thresholds. However, if the inputs are already co-partitioned on the join keys across multiple partitions (e.g. both partitioned by hash on the join keys or by range with matching boundaries), switching toCollectLeftintroduces an unnecessaryCoalescePartitionsand broadcasts data across partitions or cluster nodes.Preserving
PartitionMode::Partitionedwhen the inputs already satisfy the partitioned distribution requirements allows the join to execute locally within each partition with zero shuffles or network transfers.What changes are included in this PR?
HashJoinExec::partitioned_input_distribution_requirementsandHashJoinExec::inputs_satisfy_partitioned_requirementsindatafusion/physical-plan/src/joins/hash_join/exec.rs.JoinSelection::try_collect_left(datafusion/physical-optimizer/src/join_selection.rs), skip converting toCollectLeftif inputs already satisfy partitioned distribution requirements across >1 partition.What is the testing strategy for this PR?
New unit tests.
Are there any user-facing changes?
No.